D18 把 join 演算法選擇拆開來看,最後提一句 Sort、hash 都有 spill 路徑 ,但 spill 到底怎麼運作?
今天把三件事湊在一起說:
看起來是三個題目,其實是同一個原則的三種姿勢,都在問一句話:能不能不做完全部?
延續 D17、D18 的例子:
SELECT l_orderkey, SUM(l_extendedprice) AS revenue
FROM lineitem l JOIN orders o ON l.l_orderkey = o.o_orderkey
WHERE o.o_orderdate BETWEEN '1998-01-01' AND '1998-12-31'
GROUP BY l_orderkey
ORDER BY revenue DESC
LIMIT 100
lineitem 是六億筆的事實表,orders 幾千萬筆
SUM 群組完之後 ORDER BY 要排序,機器 heap 假設只有 8 GB,lineitem 掃出來的資料光 raw 就好幾倍, SQL 要跑起來至少三個機會:
o_orderdate BETWEEN ... 這個 filter 可以在 scan 那層就剪掉大部分 ordersSELECT l_orderkey 跟 l_extendedprice,剩下 14 個欄位不用讀ORDER BY ... LIMIT 100,其實只要維持一個大小為 100 的堆疊就夠三個機會如果一個都不用,就算演算法選得再漂亮,執行時還是要面對整批資料塞不下記憶體的問題。當那一刻真的到來,就輪到 spill 上場。
DataFusion 的 SortExec 是個 pipeline breaker,D16 講過它要看完全部資料才知道最大的是哪一筆。看完全部等於全部要進記憶體,這個假設不成立的時候怎麼辦?
答案是 Graefe 2006 整理過的 external merge sort,兩個階段:
Phase 1,Run generation:一次讀 batch 進來排序,記憶體滿了就把排好的一段 run 寫到磁碟,記憶體釋放繼續讀下一段。跑完之後磁碟上會有 N 個各自排好序的 run。
Phase 2,Multi-way merge:把 N 個 run 打開,每個 run 讀一小塊到記憶體,用一個 min-heap 挑最小的吐出來,補進來繼續挑。整個過程只需要「N 個 batch + 一個 heap」的記憶體,就能產生完整排好序的輸出。
[記憶體] [磁碟]
第一段讀進來 → sort → 寫出 →→→ Run 0 [ 排好序 ]
第二段讀進來 → sort → 寫出 →→→ Run 1 [ 排好序 ]
第三段讀進來 → sort → 寫出 →→→ Run 2 [ 排好序 ]
...
第 N 段讀進來 → sort → 寫出 →→→ Run N [ 排好序 ]
Phase 2:
Run 0 ┐
Run 1 ┼→ min-heap →→→ 完整排序輸出(串流)
Run 2 ┤
... ┘
這是 60 年代就在磁帶上做的東西,跳到 SSD 世代其實沒變太多,改變的是常數:一個 run 大小從幾百 KB 拉到幾百 MB,merge 從 5-way 拉到幾十 way,寫入放大比從 3x 降到接近 1x。演算法本身還是那個演算法,也就是為什麼今天 DataFusion 的 SortExec 走這條,跟 Postgres 的 tuplesort、Spark 的 UnsafeExternalSorter 走的其實是同一條路。
跟 MemoryPool 的介面回扣 D17:SortExec 是 spillable consumer,跟 pool grow 時被拒絕,就進 spill 路徑寫一個 run 到 tmp,reservation 釋放,pool 額度回來。整個過程 pool 不介入決策,只做批准或拒絕,spill 的方向與時機都是 SortExec 自己決定的。
排序拆完,回頭看那個 SQL 開頭lineitem 六億筆,其中 o_orderdate BETWEEN ... 可能只匹配一年的 orders,join 過濾完剩下不到一億筆;SELECT 只用兩個欄位,其他 14 個根本不需要 decode。
這些東西能不能在 scan 那一層就決定,決定了成本的量級。DataFusion 主要有三種下推:
Filter pushdown:把 WHERE 的謂詞下推到 scan。Parquet 這邊搭配 footer 的 min-max 統計,直接跳過整個 row group 不 decode;D18 的 dynamic filter 就是這條路走到極致。
Projection pushdown:只讀 SELECT 用到的欄位。Parquet 是欄式儲存,一個 row group 內每一欄各自一段 column chunk,沒選到的欄位那段 chunk 完全不去碰,I/O 直接省掉。
Limit pushdown:ORDER BY + LIMIT 這種組合,可以把「我只要前 100 筆」的資訊往下傳到 sort,甚至進到 scan。這一條下面另外講。
下推前後的差別通常不是「快兩倍」,是「掃 100 GB」變成「掃 500 MB」:
不下推 下推之後
┌─────────────────────────┐ ┌───┐
│ 讀完整 Parquet + 全欄位 │ │ ▓ │ 只讀符合 predicate
│ 全部 decode │ │ ▓ │ 的 row group +
│ │ │ ▓ │ 只選要的欄位
│ ~ 100 GB I/O + decode │ └───┘
│ │ ~ 500 MB I/O + decode
└─────────────────────────┘
每一筆都要走完整 只有留下來的資料
pipeline 需要往上走
下推的邏輯位置在 optimizer,不在 executor:optimizer 跑優化 rule 的時候把 Filter / Projection / Limit 這幾個節點「沉」到 TableScan 附近,交給 TableProvider 決定要不要真的接手(用一個「支援哪些 predicate、不支援哪些」的回傳來協商)。這條契約會在 D20 iceberg-rust 那一天正式登場,今天知道位置就好。
回到 ORDER BY revenue DESC LIMIT 100。如果按照 external sort 的定義,要先把整個群組後的資料排完序,才能拿到前 100 筆。可是想一下,真的需要排完嗎?
Top-K 用一個大小為 K 的 heap 就能拿到前 K 筆,時間複雜度 O(N log K),記憶體只要 K 那麼多,不是 O(N log N) 也不是 O(N) 記憶體。DataFusion 的 SortExec 看到下游有 LIMIT K 且 K 相對小時,會走 TopK 這條路:
原本要 sort 六億筆全排完,變成掃過去比 100 筆的 heap,六億次 O(log 100) ≈ 7 次比較。而且完全不需要 spill,因為 heap 那 100 筆一直塞得下。
這是「排序」在遇到「limit」這個資訊之後產生的演算法降級,從外部排序退回到堆疊維護。降級一詞聽起來像慢版本,其實這個「降級」是更快的版本,因為 K 遠小於 N。
回到今天的問題,排序、Pushdown、Top-K 看起來三個獨立題目,但都在做同一件事:在能夠決定這一批不需要的最早時機決定它。
三個看起來一個省 I/O、一個省 memory、一個救命,其實共用一個原則:每一份被留下來繼續往上走的資料,都應該有理由留下來
Shuffle 為什麼是 pipeline breaker 的極端形式:一般 pipeline breaker 是「這個 operator 要看完上游全部才能吐第一筆」,例如 sort、hash join build、aggregate final。Shuffle 更狠,它是跨 process、跨機器的重新分區,所有 mapper 端要把自己這一份資料按 partition 寫到磁碟,reducer 端要等到所有 mapper 都寫完才能開始讀。這個「等」比一般 breaker 多了網路加檔案系統的往返,時間量級再往上跳。所以每一個 shuffle 邊界對 planner 而言都是「非常昂貴」的邊界,能不 shuffle 就不 shuffle,能重用 shuffle 結果就重用。
Spill 的 fallback 路徑是這樣的優先順序:Comet 這邊 native sort 遇到記憶體不夠,第一選擇是自己 spill 到本地 tmp,跟 DataFusion 的路徑一樣。但 Comet 的 MemoryPool 有一條額度線是跟 Spark 的 UnifiedMemoryManager 綁定,Spark 那邊也吃緊時,Comet 這一 batch 會直接標記 unsupported,整個子樹 fallback 給 Spark 的原生 operator 處理,付一次 C2R 的代價。這條 fallback 路徑不是為了效能而設計,是為了保證即使記憶體壓力極大時查詢還能跑完。
Comet 1.0.0 有原生 shuffle 實作,走 Arrow IPC 序列化加上 native 這邊自己管 spill,跟 Spark 原生 shuffle 走完全不同的檔案格式。這條線的細節(分區、序列化、壓縮的三個決策點)留到 D29 shuffle 那一天再拆。
排序超過記憶體時走外部排序,run 生成加 multi-way merge,這是 60 年代的老演算法在 SSD 世代還沒被淘汰。Pushdown 把「能不能不讀」問到 scan 那一層,是最省的省;Top-K 把「能不能不留」問到 sort 那一層,把 O(N log N) 降到 O(N log K)。三件事看起來獨立,共用同一個原則:沒有理由留下來的資料,越早去除越好。Comet 這邊 shuffle 是 pipeline breaker 的極端形式,spill 有從 native 走到 fallback 的降級路徑保命。
明天 D20 進入第四層 iceberg-rust:一個表格式要滿足什麼契約,引擎才能真正把它當資料源?
那就明天見~
SortExec API 文件: https://docs.rs/datafusion/latest/datafusion/physical_plan/sorts/sort/struct.SortExec.html